D16 講過 DataFusion 怎麼執行(拉式、batch、tokio)、D17 講記憶體資源的協調。今天要討論是什麼東西這麼關鍵,幾乎決定整條 SQL 查詢的成本與資源用量、planner 有最多路徑可以選?
JOIN
同一句 A JOIN B ON A.x = B.x,DataFusion 至少有四種完全不同的執行策略可以走:hash、sort-merge、symmetric hash、nested loop。四個的記憶體用量、磁碟寫入、CPU 需求、能不能串流全都不一樣。
Planner 憑什麼選?表面上是一棵很短的決策樹,底下是一整套統計 + 成本模型 + 動態調整。今天把這條路徑拆開,順便看一個很優雅的技巧 dynamic filter:規劃期看的資訊不夠,執行期把 build side 建好的東西「橫向」傳過去,讓 probe side 再剪一次。
先看一個很日常的 join:
SELECT o.o_orderkey, c.c_name
FROM orders o JOIN customer c ON o.o_custkey = c.c_custkey
WHERE c.c_region = 'ASIA'
Planner 看到這句,要決定的第一件事不是「怎麼跑」,而是「這個 join 適用哪些演算法」。
DataFusion 主要選項有四個:
四個都能跑同一句 SQL,但成本差好幾個數量級,挑錯的代價不是「慢一點」,是「這條 query 跑不完」。
四個演算法裡有三個都要靠匹配 key:hash 打 bucket、sort-merge 兩邊一起走。
這些只能回答一種問題「這兩個值相等嗎」。
如果 JOIN 條件是 >、<、BETWEEN 這種範圍,或更誇張的 ON true(笛卡兒積),沒有 key 可以打,也沒有排序能一起走。這時候只剩下 nested loop。
所以第一問一定是:有沒有等值謂詞?
Nested loop 的複雜度是 O(m × n),兩邊都是百萬列的時候直接不用想。實務上有個常見加速:broadcast NLJ,小邊 broadcast 到每個 partition、大邊在 partition 內走 loop。理論複雜度沒變,但小邊留在 CPU cache 內、常數項會小很多。
SortMergeJoin 需要兩邊輸入都已排序,乍看很吃虧為了 join 額外做 sort,比直接 hash 慢。
但它有兩個場景會贏:
Spark 3.x 之後大 join 常常預設走 SortMerge 就是這個道理,事實表對事實表,兩邊都是幾億筆,hash 直接跪。
Comet 也 offload SMJ,跟 Spark 預設一致。
如果沒有天然排序,就走 hash 這條路。但 hash 有個隱含假設:build 側可以先物化完
下游要 await build 建完才能開始 probe。
大部分批次查詢都成立,可是有些場景不行:
這時候用 SymmetricHashJoin,兩邊各建一份 hash table,一邊來一 batch 就 probe 另一邊 + 插入自己這邊。代價是兩份 hash table 的記憶體,比一般 hash 貴一倍以上。
DataFusion 主要拿這個處理串流,Comet 是批次引擎,這個一般不 offload。
走到 hash 這條路,還剩最後一問哪一邊當 build、哪一邊當 probe?
答案很直白:挑估算 row count 小的當 build。build side 要整個塞進 hash table,越小記憶體越省,probe side 越大 CPU 就一次消化越多。
但這裡藏了一個很重的地雷如果估算 build side 的 row count 錯誤,會導致記憶體不足、hash table 需要 spill 到 disk,效能嚴重下降,查詢時間可能多出數倍
如果統計不準、把大表當 build,hash table 就爆到要 spill。
回想 D17 中,HashJoin 的 build 側是 spillable 的,MemoryPool 拒絕 reservation 時它會把部分 hash bucket 寫到 disk。
Spill 之後 probe 進來每一筆都要先看記憶體、再看 disk,兩趟 I/O 疊上去。
所以你會看到 CBO(cost-based optimizer)這個詞在 join 場景特別重要——選錯 build side 的代價是幾倍,不是幾成。
把上面拼起來:
┌──────────────────────────┐
│ 有等值謂詞嗎? │
└────┬──────────────┬──────┘
│沒有 │有
▼ ▼
NestedLoopJoin 兩邊都排序好?
(或 broadcast) ┌────┴────┐
│是 │否
▼ ▼
SortMergeJoin 兩邊都不能先物化?
┌───┴───┐
│是 │否
▼ ▼
SymmetricHashJoin HashJoin
(挑小的當 build)
實務上還有更多細節:left / right / full outer 各有偏好、broadcast vs shuffle 是另一個維度、null-aware anti join 有特殊處理。但主幹就這條。
上面決策樹都是 planner 在規劃期做的事,看 SQL、看統計、選演算法
可是有個問題:很多資訊要跑起來才知道。
回到開頭那個例子:orders JOIN customer ON o.o_custkey = c.c_custkey WHERE c.c_region = 'ASIA'。
規劃期只知道「customer 這邊會被 filter」,但 filter 完剩多少 key、剩下哪些 key,都得跑起來才知道。假設 c_region = 'ASIA' 篩完只剩 200 個 c_custkey——那 orders 本來要掃 6 億筆,其實只要看那 200 個 c_custkey 出現在哪些 row group 裡就好。
Dynamic Filter 就是這件事:
這個技巧的正式名字是 sideways information passing——資訊不是走 pipeline 上下游(那是 top-down),而是橫向從一個 operator 傳到另一個。
為什麼把 join filter 動態下推到 scan 是有效的?
粗糙版答案:「因為 scan 少讀了」但這個答案沒觸到重點,如果 scan 只是把資料讀進來後在 filter 節點過濾,那 dynamic filter 只是省了一段 in-memory 處理,不會有數量級差別。
真正有效的原因是 Parquet 的檔案內剪枝機制:
也就是說 dynamic filter 有效不是因為它是個 filter,而是因為它讓 scan 少做 I/O 跟 decode。如果 probe side 是 CSV(沒有檔案內剪枝機制),這個技巧退化成只是 in-memory filter,效果差很多。
還有一個時機關鍵:dynamic filter 必須在 probe scan 還沒開始(或早期)時送到。太晚 scan 已經讀完就沒意義。DataFusion 用 async channel + repartition 邊界處理這件事,也是 pull + async 執行模型的一個副作用——channel 天然可以做橫向資訊傳遞。
Comet 目前的 join 支援大致是(細節看 release note 跟 compatibility matrix):
Fallback 的代價回到 D15 接點五——整個 join 子樹回到 JVM 走 Spark 原生,這一批要從 Arrow C Data Interface 轉回 UnsafeRow(付一次 C2R)、join 完再考慮要不要轉回 Arrow 繼續 native(付一次 R2C)。所以 fallback 不是「這個 join 慢一點」,是這一批資料要多兩次跨界。
Join 演算法選擇的主幹是短短一棵決策樹:有等值 → 有序 → 能物化 → 挑小邊。挑錯的代價不是慢,是查詢跑不完。Dynamic Filter 是規劃期的決策不夠時、執行期補一刀的技巧,有效性靠 Parquet 檔案內剪枝背書——是儲存跟執行的一次成功協作。
明天 D19 接執行細節的下一題:資料超過記憶體時排序怎麼繼續?外部排序是所有 pipeline breaker 都會遇到的降級路徑,順便把「下推」跟「提前終止」這幾個省 I/O 的手法一起說完。
那就明天見~